MBL
Go / avito / Тестовое задание / Transactional outbox и DLQ на практике
Go сложный

Transactional outbox и DLQ на практике

avitokafkapostgresreliability

Надёжная доставка между двумя сервисами через брокер — тема, где почти любой ответ «в общих словах» звучит одинаково у всех кандидатов. Разница видна только на конкретном коде: что происходит в момент сбоя, шаг за шагом.

Проблема без outbox

Наивная реализация: закоммитить транзакцию создания заказа, потом опубликовать сообщение в Kafka.

tx.Commit(ctx)                  // заказ уже в БД
producer.Publish(ctx, msg)      // ...а тут процесс может упасть

Между этими двумя строками — окно, в котором заказ существует в БД, но establishment-service никогда не узнает о нём. Переставить местами (сначала publish, потом commit) не решает проблему, а только переносит её: теперь можно опубликовать сообщение о заказе, которого в БД в итоге не оказалось (транзакция откатилась).

Решение: запись в outbox той же транзакцией

err := txManager.WithinTx(ctx, func(ctx context.Context) error {
    order, err := orderRepo.Create(ctx, order)
    if err != nil {
        return err
    }
    payload, _ := json.Marshal(toOrderCreatedMessage(order))
    return outboxWriter.Enqueue(ctx, domain.TopicOrdersNew, order.ID.String(), payload)
})

Enqueue пишет строку в outbox_events через тот же q(ctx, pool), что и orderRepo.Create — обе операции живут в одной транзакции Postgres. Либо обе применяются, либо ни одна. Публикация в Kafka здесь вообще не участвует — она происходит позже, отдельным процессом.

Фоновый relay

func (r *Relay) tick(ctx context.Context) {
    rows, _ := r.store.FetchUnpublished(ctx, r.batchSize)
    for _, row := range rows {
        if err := r.producer.Publish(ctx, row.Topic, row.Key, row.Payload); err != nil {
            continue // ретрай следующим тиком, published_at не трогаем
        }
        if hook, ok := r.onPublished[row.Topic]; ok {
            if err := hook(ctx, row); err != nil {
                continue // тот же принцип: не пометили — хук вызовется снова
            }
        }
        r.store.MarkPublished(ctx, row.ID)
    }
}

Три исхода на каждую строку:

  • публикация упала → пропустить, следующий тик повторит с нуля;
  • публикация прошла, но пост-хук (например, MarkSentToEstablishment) упал → тоже пропустить, не помечать published_at — иначе хук никогда не будет вызван повторно, а строка уже будет выглядеть «обработанной»;
  • оба шага успешны → пометить.

Ключевое решение — пост-хук выполняется после publish, и его сбой откатывает не саму публикацию (она уже произошла в Kafka), а только отметку published_at. Значит одно и то же сообщение может уйти в Kafka дважды, если хук падал между попытками — это осознанно допустимо, потому что consumer на другой стороне идемпотентен (см. «Пять разных идемпотентностей»).

Юнит-тест на этот цикл — не e2e, а таблично, с фейковыми Store/Producer, покрывает ровно эти три исхода плюс топик без зарегистрированного хука:

func TestRelay_Tick_PublishFails(t *testing.T) {
    store := &fakeStore{rows: []outbox.Row{{ID: 1, Topic: "orders.new"}}}
    producer := &fakeProducer{err: errors.New("kafka down")}
    r := outbox.NewRelay(store, producer, nil, 10)

    r.Tick(context.Background())

    if store.markedPublished[1] {
        t.Fatal("must not mark published when publish failed")
    }
}

DLQ: offset — это watermark, не индивидуальный ack

Ключевая вещь, которую легко сказать неправильно на собесе: Kafka-consumer не может «отложить одно сообщение и обработать следующее», как это можно сделать, например, с SQS-очередью с индивидуальными ack по сообщению.

Offset — это одно число на партицию: «всё до этой позиции обработано». Закоммитить offset сообщения N+1, оставив N необработанным, невозможно осмысленно — коммит N+1 неявно означает «N тоже подтверждено».

Отсюда — обязательное разделение ошибок на два класса до решения, коммитить offset или нет:

func process(ctx context.Context, msg kafka.Message, handler HandlerFunc) {
    isDomainErr, err := handler(ctx, msg.Key, msg.Value)
    if err == nil {
        consumer.CommitMessages(ctx, msg)
        return
    }
    if isDomainErr {
        toDLQ(ctx, msg, err)   // никогда не станет валидным само по себе
        return                  // toDLQ коммитит offset сам, после успешной публикации в DLQ
    }
    // транзиентная: до 5 попыток с backoff, ПОТОМ тоже в DLQ
    retryWithBackoff(ctx, msg, handler)
}

Доменная ошибка (ErrInvalidTransition — переход не разрешён и не совпадает с текущим статусом) — ретраить бессмысленно, сообщение никогда не станет валидным сколько его ни пытайся обработать снова. Сразу в DLQ-топик (orders.status_updated.dlq), затем commit основного offset.

Транзиентная ошибка (БД временно недоступна) — до 5 попыток с экспоненциальным backoff внутри обработки этого же сообщения, не через повторный FetchMessage. Если все попытки исчерпаны — тоже в DLQ, а не бесконечный ретрай. Без этого одно «вечно падающее» сообщение блокирует партицию для всех последующих заказов навсегда.

func toDLQ(ctx context.Context, msg kafka.Message, cause error) {
    dlqPayload := wrapWithReason(msg.Value, cause)
    if err := dlqProducer.Publish(ctx, dlqTopic, msg.Key, dlqPayload); err != nil {
        return // публикация в DLQ не удалась — НЕ коммитим offset основного топика,
                // сообщение переобработается при рестарте, а не потеряется
    }
    consumer.CommitMessages(ctx, msg)
}

Даже путь в DLQ не коммитит offset, пока публикация в сам DLQ-топик не подтверждена — иначе сбой Kafka именно в этот момент тихо терял бы сообщение вместо того, чтобы просто повторить попытку позже.

Проверено e2e, не только вручную

e2e/dlq_test.go публикует напрямую в orders.status_updated с хоста (через EXTERNAL://localhost:29092 listener, специально добавленный в docker-compose.yml для этого теста) невалидный переход sent_to_establishment → completed, минуя establishment-service:

func TestOrderStatusUpdated_InvalidTransitionGoesToDLQ(t *testing.T) {
    // ... публикация невалидного перехода напрямую в Kafka ...

    assertMessageInDLQ(t, dlqReader, orderID) // читает с kafka.FirstOffset —
        // свежий consumer-group иначе стартовал бы с "сейчас" и пропустил
        // уже опубликованное сообщение

    assertOrderStatusUnchanged(t, orderID, "sent_to_establishment")

    publishValidTransition(t, orderID, "accepted") // тот же ключ = та же партиция
    waitForStatus(t, orderID, "accepted")           // партиция не заблокирована
}

Три вещи проверяются вместе: сообщение действительно попало в DLQ, заказ не сдвинулся с места, и — самое важное — партиция не заблокирована poison-message'ем: следующее валидное сообщение на том же ключе (значит на той же партиции) обрабатывается штатно.

Запусти сам: как ведёт себя relay на трёх исходах

Логика tick (публикация упала / хук упал / всё прошло) не требует ни Kafka, ни Postgres, чтобы её увидеть — вот та же машина состояний на fake-реализациях, целиком:

на play.golang.org

Вывод показывает ровно то, что описано выше прозой: A не опубликован (publish упал, ретрай следующим тиком), B фактически ушёл в Kafka, но не помечен published (хук упал — при следующем тике Publish вызовется ещё раз, это осознанно допустимо), только C дошёл до конца цикла.

Самопроверка 0 / 4
Могу объяснить, какую именно потерю сообщений закрывает outbox
Понимаю, почему relay не помечает published_at при сбое пост-хука
Знаю, почему нельзя просто закоммитить offset N+1 и ретраить N отдельно
Различаю доменную ошибку (сразу в DLQ) и транзиентную (ретраить)
Как усвоено?